Mark system Nexus envelope payloads#1667
Conversation
4cf4eaa to
b489b0e
Compare
| async def _visit_nexus_operation_input_payload( | ||
| self, | ||
| fs: VisitorFunctions, | ||
| endpoint: str, | ||
| payload: Payload, | ||
| ) -> None: | ||
| new_payload = await temporalio.nexus.system._maybe_visit_payload( | ||
| endpoint, | ||
| payload, | ||
| await self._visit_temporal_api_common_v1_Payload(fs, payload) | ||
|
|
||
| async def _visit_temporal_api_common_v1_Payload( | ||
| self, fs: VisitorFunctions, o: Payload | ||
| ): | ||
| new_payload = await temporalio.nexus.system.maybe_visit_payload( | ||
| o, | ||
| fs, | ||
| self.skip_search_attributes, | ||
| ) | ||
| if new_payload is None: | ||
| await self._visit_temporal_api_common_v1_Payload(fs, payload) | ||
| await fs.visit_payload(o) | ||
| return | ||
|
|
||
| if new_payload is not payload: | ||
| payload.CopyFrom(new_payload) | ||
| await fs.visit_system_nexus_envelope(payload) | ||
|
|
||
| async def _visit_temporal_api_common_v1_Payload( | ||
| self, fs: VisitorFunctions, o: Payload | ||
| ): | ||
| await fs.visit_payload(o) | ||
| if new_payload is not o: | ||
| o.CopyFrom(new_payload) | ||
| await fs.visit_system_nexus_envelope(o) |
There was a problem hiding this comment.
Not sure if the formatter allows it, but I'd put put some of these args and params on the same line and rename o to payload
There was a problem hiding this comment.
It would be a little painful I think because it is all generated code, and is using some of the same names as all other functions.
There was a problem hiding this comment.
I take it back, looks like this one isn't part of that, so it should be easy
| ) | ||
| visitor = SystemNexusVisitor() | ||
|
|
||
| await PayloadVisitor().visit(visitor, comp) |
There was a problem hiding this comment.
Codex observed that this fails if concurrency_limit>1.
There was a problem hiding this comment.
Can you elaborate? And do you think it matters for some reason?
There was a problem hiding this comment.
I mean the test fails if line 253 were replaced with:
await PayloadVisitor(concurrency_limit=3).visit(visitor, comp)It sounded to me like an issue with the implementation, but I didn't understand well enough so I thought I'd mention it. If you don't think it matters then it probably doesn't matter
What changed
System Nexus envelope payloads encoded by the system Nexus payload converter now carry a
__temporal_system_payloadmetadata marker with valueb"true".Payload visitors now use that marker to detect system Nexus envelopes whenever they encounter an already-encoded payload. This lets nested user payloads inside the envelope still be decoded and visited even when the payload is no longer accompanied by a Nexus endpoint field.
Detection behavior
Detection is metadata-only. Unmarked payloads are treated as regular payloads, even if they are used as the input for the Temporal system Nexus endpoint. Marked payloads are treated as system Nexus envelopes regardless of endpoint.
Tests
./.venv/bin/pytest tests/nexus/test_temporal_system_nexus.py -q./.venv/bin/pytest tests/worker/test_visitor.py -qpoe lint